[Deferred] Dead letter queue for deferred tasks - #383
Conversation
1e367e1 to
82e1ab4
Compare
… dead task into DLQ after all retries got exhausted
82e1ab4 to
e144296
Compare
Deferred Backends ApproachClass dependenciesclassDiagram
class Queue {
-backend
+initialize(backend)
+enqueue(...)
+schedule(...)
}
class Nil {
+TasksBackend tasks
+DeadTasksBackend dead_tasks
+initialize(**)
}
class Disk {
+TasksBackend tasks
+DeadTasksBackend dead_tasks
+initialize(tasks_options:, dead_tasks_options:)
}
class NilTasksBackend["Nil::TasksBackend"] {
+add(...)
+remove(...)
+pending_tasks()
}
class NilDeadTasksBackend["Nil::DeadTasksBackend"] {
+add(...)
+remove(...)
+retry(...)
}
class DiskTasksBackend["Disk::TasksBackend"] {
- path
- prefix
- fsync_frequency
+add(context, publish_at:, task_id:)
+remove(task_id)
+pending_tasks()
}
class DiskDeadTasksBackend["Disk::DeadTasksBackend"] {
- path
- prefix
+add(context, exception:, task_id:)
+remove(task_id)
+retry(...)
}
Queue --> Nil : backend
Queue --> Disk : backend
Nil *-- NilTasksBackend : tasks
Nil *-- NilDeadTasksBackend : dead_tasks
Disk *-- DiskTasksBackend : tasks
Disk *-- DiskDeadTasksBackend : dead_tasks
Runtime call shape: backend.tasks.add(...)
backend.tasks.remove(...)
backend.tasks.pending_tasks
backend.dead_tasks.add(...)
backend.dead_tasks.remove(...)File interfaces
class Rage::Deferred::Backends::Disk
attr_reader :tasks, :dead_tasks
def initialize(tasks_options: {}, dead_tasks_options: {})
@tasks = TasksBackend.new(**tasks_options)
@dead_tasks = DeadTasksBackend.new(**dead_tasks_options)
end
class TasksBackend
def initialize(path:, prefix:, fsync_frequency:)
# WAL setup (current Disk body)
end
def add(context, publish_at: nil, task_id: nil) end
def remove(task_id) end
def pending_tasks end
end
class DeadTasksBackend
def initialize(path:, prefix:)
...
end
def add(context, exception:, task_id:) end
def remove(task_id) end
def retry(...) end
end
end
class Rage::Deferred::Backends::Nil
attr_reader :tasks, :dead_tasks
def initialize(**)
@tasks = TasksBackend.new
@dead_tasks = DeadTasksBackend.new
end
class TasksBackend
def add(_, **) end
def remove(_) end
def pending_tasks = []
end
class DeadTasksBackend
def add(_, **) end
def remove(_) end
def retry(_, **) end
end
endConfiguration
Rage.configure do
config.deferred.backend = :disk
config.deferred.tasks = { path: "storage", prefix: "deferred-", fsync_frequency: 0.5 }
config.deferred.dead_tasks = { path: "storage", prefix: "dead_tasks-" }
end
@backend_class.new(tasks_options: @tasks_options, dead_tasks_options: @dead_tasks_options)For |
|
@rsamoilov @serhii-sadovskyi please take a look at the proposed approach above for the deferred backends implementation. I'd like to hear your thoughts! |
|
Hey @alex-rogachev , This looks great. Love the diagram! Several thoughts:
This would be a breaking change, which we would ideally want to avoid.
This exposes the fact that dead tasks are stored in another file, which is an internal implementation detail that shouldn't be exposed. Consider a hypothetical Redis backend - would it need a separate The DLQ backend should be reusing the same
If this method is supposed to reenqueue the task and remove it from the DLQ, then it should be part of another class, potentially |
All these points make sense. We'll consider them if we select this approach! |
Work In Progress